internal/api
import "github.com/neochaotic/leoflow/internal/api"
Package api implements the Airflow-compatible HTTP control plane (ADR 0007).
Index
- Constants
- Variables
- func AbortProblem(c *gin.Context, status int, title, detail string)
- func CORS(allowed []string) gin.HandlerFunc
- func DevBypassAuth() gin.HandlerFunc
- func JWTAuth(authn auth.Authenticator) gin.HandlerFunc
- func NewServer(deps Dependencies) *gin.Engine
- func NoStoreOnVolatileRoutes() gin.HandlerFunc
- func ObservabilityHandler(registry *prometheus.Registry, checks map[string]HealthChecker) http.Handler
- func Observe(metrics Metrics, tracer trace.Tracer) gin.HandlerFunc
- func RequestID() gin.HandlerFunc
- func RequirePermission(action, resource string) gin.HandlerFunc
- func StructuredLogger(logger *slog.Logger) gin.HandlerFunc
- func UserFromContext(c *gin.Context) (*auth.User, bool)
- type AuditLogReader
- type AuditWriter
- type AuthAuditWriter
- type ConnectionStore
- type ConnectionTester
- type DagLatestRunsReader
- type DagRepository
- type DagRunRepository
- type DagSpecReader
- type DagVersionLister
- type DagVersionRepository
- type DashboardStatsReader
- type Dependencies
- type ExecutorInfo
- type FavoriteStore
- type HealthChecker
- type Heartbeater
- type ImportErrorStore
- type LogReader
- type Metrics
- type OIDCUserStore
- type PoolStore
- type Problem
- type TaskInstanceRepository
- type TaskSummaryReader
- type TokenRenewer
- type UIServer
- type UserAuditWriter
- type UserStore
- type VariableStore
- type WorkspaceFS
- type XComReader
Constants
DefaultUIAutoRefreshIntervalSeconds is the production-safe value returned by /ui/config when no explicit override is configured. Lite overrides this to a smaller value (typically 5s) for a snappy inner-loop dev experience; Pro keeps 30s so the SPA’s polling does not hammer a shared metadata DB.
const DefaultUIAutoRefreshIntervalSeconds = 30
Variables
ErrNotFound is returned by repositories when a resource does not exist.
var ErrNotFound = domain.ErrNotFound
func AbortProblem
func AbortProblem(c *gin.Context, status int, title, detail string)
AbortProblem writes an RFC 7807 problem response and stops the handler chain.
func CORS
func CORS(allowed []string) gin.HandlerFunc
CORS allows the configured origins (use “*” to allow any).
func DevBypassAuth
func DevBypassAuth() gin.HandlerFunc
DevBypassAuth authenticates EVERY request as a fixed admin user, with no token required. It exists solely for `leoflow dev` (the local, unsandboxed loop) so a developer reaches the UI without logging in. It must only be wired under the explicit dev opt-in (config auth.dev_no_auth); the server logs a prominent warning when it is active. NEVER enable this in production.
func JWTAuth
func JWTAuth(authn auth.Authenticator) gin.HandlerFunc
JWTAuth validates the bearer token on protected routes and stores the user.
func NewServer
func NewServer(deps Dependencies) *gin.Engine
NewServer builds the gin engine with the full middleware chain, health and metrics endpoints, embedded Scalar docs, and the auth token endpoint.
func NoStoreOnVolatileRoutes
func NoStoreOnVolatileRoutes() gin.HandlerFunc
NoStoreOnVolatileRoutes stamps every response served from the SPA-facing JSON surface (`/api/v2/*` and `/ui/*`) with `Cache-Control: no-store, must-revalidate` so the browser HTTP cache does not return a pre-mutation payload after a PATCH/POST/DELETE.
Why this exists (#211, #271): mark-state PATCH succeeds in single-digit ms; TanStack Query then invalidates its in-memory cache and re-fetches. Without this header, the browser’s HTTP cache layer can serve the OLD response to that re-fetch (the original GET response had no explicit caching directive, so the browser falls back to heuristic caching). The SPA renders stale state until the next “natural” refresh — the observable symptom is “marcar como falha demora uma eternidade”.
Static assets (`/ide/vs/*` for the Monaco bundle) are content-hashed and SHOULD cache, so they are explicitly excluded.
We deliberately use “no-store” rather than “no-cache”: no-store forbids the browser from writing the response anywhere, which is the strongest guarantee we can give a TanStack-backed SPA. “must-revalidate” is added for older intermediaries (proxies / SW) that may not honor no-store alone. This is ADR-0017-compatible: no SPA changes.
func ObservabilityHandler
func ObservabilityHandler(registry *prometheus.Registry, checks map[string]HealthChecker) http.Handler
ObservabilityHandler builds the handler served on the metrics listener: the Prometheus /metrics endpoint plus the same /healthz and /readyz the full API exposes (reusing the very same handlers, so the semantics are identical — trivial liveness, dependency-pinging readiness).
Roles that do not serve the full API — the ADR 0049 scheduler role — still need a liveness/readiness surface for the kubelet’s probes. The metrics listener runs in every role, so mounting health here gives a scheduler-only pod a probe target without exposing the API, auth, or UI. The full API keeps its own /healthz and /readyz on the HTTP port, so the “all” role is unchanged; this is purely additive on the metrics port.
This handler is intentionally unauthenticated (probes carry no token) and does not run the API middleware chain; it is the same trust level as scraping /metrics, which is already public.
func Observe
func Observe(metrics Metrics, tracer trace.Tracer) gin.HandlerFunc
Observe wraps each request in an OTel span and records HTTP metrics (ADR 0010). A nil tracer falls back to the global (no-op) tracer; nil metrics are skipped, so the middleware is safe in tests.
func RequestID
func RequestID() gin.HandlerFunc
RequestID assigns a request id (honoring an inbound X-Request-Id) and echoes it.
func RequirePermission
func RequirePermission(action, resource string) gin.HandlerFunc
RequirePermission enforces an RBAC permission on a route.
func StructuredLogger
func StructuredLogger(logger *slog.Logger) gin.HandlerFunc
StructuredLogger logs one structured line per request (ADR 0010).
func UserFromContext
func UserFromContext(c *gin.Context) (*auth.User, bool)
UserFromContext returns the authenticated user stored by JWTAuth.
type AuditLogReader
AuditLogReader lists recorded actions for the Audit Log view. dagID == "" means no DAG filter.
type AuditLogReader interface {
ListAuditLogs(ctx context.Context, tenant, dagID string, limit, offset int) ([]domain.AuditLogEntry, int, error)
}
type AuditWriter
AuditWriter records task-level actions (clear, mark state) for the Audit Log view, with the acting user and the run/task in the entry.
type AuditWriter interface {
RecordTaskActionAudit(ctx context.Context, tenant, userID, action, dagID, runID, taskID string, tryNumber int) error
}
type AuthAuditWriter
AuthAuditWriter records authentication events to the audit sink (H5). It is best-effort: a write error never changes the auth outcome.
type AuthAuditWriter interface {
RecordAuthEvent(ctx context.Context, tenant, actorUserID, action, email, outcome string, extra map[string]string) error
}
type ConnectionStore
ConnectionStore reads and writes Airflow-style Connections for the Admin UI. Password is encrypted at rest by the store (ADR 0019) and never returned.
type ConnectionStore interface {
ListConnections(ctx context.Context, tenant string, limit, offset int) ([]domain.Connection, int, error)
GetConnection(ctx context.Context, tenant, connID string) (domain.Connection, error)
SetConnection(ctx context.Context, tenant string, c domain.Connection) error
DeleteConnection(ctx context.Context, tenant, connID string) error
}
type ConnectionTester
ConnectionTester checks whether a connection is well-formed. The default implementation validates STRUCTURE only and makes no network call: the Go control plane must never reach out to a user-configured host (SSRF / internal port-scan — go/request-forgery), and “reachable from the control plane” is the wrong question anyway, since a connection is used in the task’s network scope, not the control plane’s. Live reachability/auth is tested where the connection is actually used (the task/executor) — tracked as a follow-up.
type ConnectionTester interface {
Test(ctx context.Context, c domain.Connection) (ok bool, message string)
}
type DagLatestRunsReader
DagLatestRunsReader fetches the most-recent runs for a set of DAGs in one query, so /ui/dags can embed run history without an N+1.
type DagLatestRunsReader interface {
LatestRunsForDags(ctx context.Context, tenant string, dagIDs []string, perDag int) (map[string][]domain.DagRun, error)
}
type DagRepository
DagRepository reads, updates, and deletes registered DAGs.
type DagRepository interface {
ListDags(ctx context.Context, tenant string, limit, offset int) ([]domain.DAG, int, error)
GetDag(ctx context.Context, tenant, dagID string) (domain.DAG, error)
SetPaused(ctx context.Context, tenant, dagID string, paused bool) (domain.DAG, error)
DeleteDag(ctx context.Context, tenant, dagID string) error
ClearDagHistory(ctx context.Context, tenant, dagID string) error
ListDagsFiltered(ctx context.Context, tenant, runState string, paused *bool, limit, offset int) ([]domain.DAG, int, error)
}
type DagRunRepository
DagRunRepository reads and creates DAG runs.
type DagRunRepository interface {
ListDagRuns(ctx context.Context, tenant, dagID string, limit, offset int) ([]domain.DagRun, int, error)
GetDagRun(ctx context.Context, tenant, dagID, runID string) (domain.DagRun, error)
CreateDagRun(ctx context.Context, tenant, dagID string, run domain.DagRun) (domain.DagRun, error)
SetDagRunState(ctx context.Context, tenant, dagID, runID, state string) error
DeleteDagRun(ctx context.Context, tenant, dagID, runID string) error
}
type DagSpecReader
DagSpecReader reads the parsed spec of a DAG’s current version, the source of task topology for the grid and graph views.
type DagSpecReader interface {
GetCurrentSpec(ctx context.Context, tenant, dagID string) (domain.DAGSpec, error)
}
type DagVersionLister
DagVersionLister lists a DAG’s registered versions. The Airflow UI fetches this to resolve a version_number before requesting version-scoped structure (the Graph view); without it the graph never loads. See docs/ui-compatibility.md.
type DagVersionLister interface {
ListDagVersions(ctx context.Context, tenant, dagID string) ([]domain.DagVersion, error)
}
type DagVersionRepository
DagVersionRepository registers compiled DAG versions.
type DagVersionRepository interface {
// RegisterDagVersion upserts the DAG and inserts a version keyed by
// specHash, reporting whether a new version was created (false if the hash
// already existed — the push is idempotent).
RegisterDagVersion(ctx context.Context, tenant string, spec domain.DAGSpec, specHash string) (bool, error)
}
type DashboardStatsReader
DashboardStatsReader backs the home dashboard widgets with real counts.
type DashboardStatsReader interface {
DagStats(ctx context.Context, tenant string) (domain.DagStats, error)
HistoricalMetrics(ctx context.Context, tenant string, since, until time.Time) (domain.HistoricalMetrics, error)
}
type Dependencies
Dependencies bundles everything the HTTP server needs.
type Dependencies struct {
Logger *slog.Logger
Authenticator auth.Authenticator
RateLimiter *auth.RateLimiter
Registry *prometheus.Registry
Metrics Metrics
Tracer trace.Tracer
HealthChecks map[string]HealthChecker
CORSOrigins []string
// TrustedProxies is the set of proxy IPs/CIDRs whose X-Forwarded-For header
// gin will honor when resolving c.ClientIP(). Empty/nil trusts NO proxy, so
// ClientIP is the direct peer and a spoofed XFF cannot forge the client IP
// (audit H1). A Pro deployment behind an ingress sets this to the ingress
// CIDR so per-client rate-limiting and audit see the real client.
TrustedProxies []string
TokenTTLSecs int
// TokenRenewer re-mints a still-valid user bearer with a fresh short TTL so a
// long CLI/dev session need not re-login every TokenTTLSecs (aresta #5). Nil
// leaves the renew route unregistered (renewal simply unavailable). In practice
// it is the same *auth.JWTAuthenticator as Authenticator.
TokenRenewer TokenRenewer
// TokenMaxLifetimeSecs is the hard ceiling on a renewed session's total age
// since first login; past it, renewal is refused and the user must
// re-authenticate. Non-positive disables the ceiling.
TokenMaxLifetimeSecs int
// InstanceName is shown in the UI navbar (Airflow's instance_name). Empty
// falls back to "Leoflow"; `leoflow dev` sets it to mark the DEV environment.
InstanceName string
// UIAutoRefreshIntervalSeconds controls the SPA's polling cadence for DAG /
// DagRun / task-instance state refresh (Airflow's auto_refresh_interval).
// Non-positive (the zero default) falls back to DefaultUIAutoRefreshIntervalSeconds
// (30s, production-safe). `leoflow lite` sets it to ~5s for a snappy inner loop.
UIAutoRefreshIntervalSeconds int
// DevNoAuth replaces JWT auth with a dev-only bypass that authenticates every
// request as an admin (no login). It is for `leoflow dev` only and must never
// be set in production. See DevBypassAuth.
DevNoAuth bool
// Edition marks the running edition ("pro", "lite", or empty). It gates
// Pro-only surfaces: named-pool CRUD is registered as real endpoints only when
// Edition == "pro" (ADR 0053), otherwise the Pools screen gets the graceful
// empty-collection stub, matching how the scheduler's pool gate is Pro-gated.
Edition string
// Resource repositories. Routes for nil repositories are not registered.
Dags DagRepository
DagRuns DagRunRepository
Tasks TaskInstanceRepository
Versions DagVersionRepository
Xcoms XComReader
Logs LogReader
Specs DagSpecReader
LatestRuns DagLatestRunsReader
TaskSummary TaskSummaryReader
DagVersions DagVersionLister
DashboardStats DashboardStatsReader
AuditLog AuditLogReader
Variables VariableStore
Users UserStore
UserAudit UserAuditWriter
Connections ConnectionStore
ConnectionTest ConnectionTester
Pools PoolStore
Favorites FavoriteStore
ImportErrors ImportErrorStore
Audit AuditWriter
ExecutorInfo ExecutorInfo
// Workspace backs the Lite web editor (ADR 0025). When nil the editor's
// filesystem API is not registered (Production, or Lite without a workspace).
Workspace WorkspaceFS
// MonacoDir is the directory holding the pinned Monaco bundle that
// `leoflow setup` fetched; the editor page is served Monaco from it. Empty or
// missing makes the page show a setup hint instead of a broken editor.
MonacoDir string
// ExamplesFS backs the IDE's "Download examples" button — typically the
// `embed.FS` shipped from the leoflow root package. Nil disables the button.
ExamplesFS fs.FS
// SchedulerHealth reports the scheduler's heartbeat for /monitor/health.
// When nil the component reports healthy (single-process role assumption).
SchedulerHealth Heartbeater
// UI serves the embedded SPA. When nil the server is API-only.
UI UIServer
// OIDC wiring (registered only when OIDCFlow is non-nil — provider: oidc).
// The JWT authenticator above stays the request-path verifier in both modes.
//
// OIDCFlow is the discovered Authorization Code + PKCE flow; nil in JWT mode,
// in which case the /api/v2/auth/oidc/* routes are not registered.
OIDCFlow *oidc.Flow
// OIDCSettings carries the role mappings, JIT policy, default_role, and
// break-glass allowlist the login flow and the credential gate read.
OIDCSettings config.OIDCSection
// OIDCUsers resolves and JIT-provisions OIDC identities (the storage repo).
OIDCUsers OIDCUserStore
// AuthAudit records authentication events (login, tenant-pin rejection, JIT,
// break-glass, logout) to the audit sink.
AuthAudit AuthAuditWriter
// JWTSecret is the HS256 secret the OIDC callback mints the app's _token with.
JWTSecret string
}
type ExecutorInfo
ExecutorInfo describes the control plane’s execution capacity. It surfaces whether pod dispatch is available — the cluster-level answer to “why is a task stuck queued” (#46/#47). The stock Airflow UI has no widget for it, but operators (curl/monitoring) and a future custom Cluster Activity view consume it. Cluster Activity in Airflow 3.2 is otherwise the Home dashboard, already backed by /api/v2/monitor/health (#33) and /ui/dashboard/* (#39).
type ExecutorInfo struct {
PodDispatchEnabled bool
TaskNamespace string
AgentControlPlaneAddr string
}
type FavoriteStore
FavoriteStore persists per-user DAG favorites (the DAG-list star).
type FavoriteStore interface {
AddFavorite(ctx context.Context, tenant, userID, dagID string) error
RemoveFavorite(ctx context.Context, tenant, userID, dagID string) error
FavoriteDagIDs(ctx context.Context, tenant, userID string) (map[string]bool, error)
}
type HealthChecker
HealthChecker reports dependency health for readiness checks.
type HealthChecker interface {
Ping(ctx context.Context) error
}
type Heartbeater
Heartbeater reports a long-running component’s liveness and last heartbeat for the monitor health endpoint. The scheduler implements it.
type Heartbeater interface {
Heartbeat() (healthy bool, last time.Time)
}
type ImportErrorStore
ImportErrorStore reads and writes DAG parse/compile errors that back Airflow’s “Import Errors” banner on the home dashboard. The `leoflow dev` watcher writes an entry on a failed compile and clears it on the next good compile; the public GET /api/v2/importErrors feed is what the UI polls.
type ImportErrorStore interface {
ListImportErrors(ctx context.Context, tenant string) ([]domain.ImportError, error)
SetImportError(ctx context.Context, tenant string, e domain.ImportError) error
ClearImportError(ctx context.Context, tenant, filename string) error
}
type LogReader
LogReader streams a task attempt’s stored logs and, for running tasks, tails new lines live.
type LogReader interface {
ReadLogs(ctx context.Context, tenant, dagID, runID, taskID string, tryNumber int) (io.ReadCloser, error)
Tail(ctx context.Context, tenant, dagID, runID, taskID string, tryNumber int) (<-chan string, func(), error)
}
type Metrics
Metrics records HTTP request metrics. observability.Metrics implements it.
type Metrics interface {
RecordHTTPRequest(method, path string, status int, dur time.Duration)
}
type OIDCUserStore
OIDCUserStore resolves and just-in-time-provisions OIDC identities. storage implements it. The interface lives with its consumer (the callback handler).
type OIDCUserStore interface {
// FindUserByOIDCSubject resolves a returning identity by its (provider,
// subject) pair, returning the user, whether it is active, and
// auth.ErrUserNotFound when no row matches.
FindUserByOIDCSubject(ctx context.Context, provider, subject string) (*auth.User, bool, error)
// CreateOIDCUser JIT-provisions an OIDC-only user with the given roles.
CreateOIDCUser(ctx context.Context, tenant, email, provider, subject string, roles []string) (*auth.User, error)
// RoleExists reports whether a role name exists for the tenant, used to fail
// closed on a misconfigured default_role or role mapping.
RoleExists(ctx context.Context, tenant, role string) (bool, error)
// ReconcileUserRoles sets the user's DB roles to EXACTLY roleNames (the IdP is
// authoritative). It fails closed on a name that is not a role in the tenant,
// leaving the prior grants intact.
ReconcileUserRoles(ctx context.Context, userID string, roleNames []string) error
}
type PoolStore
PoolStore reads and writes named task pools for the Admin UI (ADR 0053 Stage 3). Pools are tenant-scoped; PoolSlotUsage reports per-pool occupancy for the Airflow slot fields.
type PoolStore interface {
ListPools(ctx context.Context, tenant string, limit, offset int) ([]domain.Pool, int, error)
GetPool(ctx context.Context, tenant, name string) (domain.Pool, error)
SetPool(ctx context.Context, tenant string, p domain.Pool) error
DeletePool(ctx context.Context, tenant, name string) error
PoolSlotUsage(ctx context.Context, tenant string) (map[string]domain.PoolUsage, error)
}
type Problem
Problem is an RFC 7807 problem-details response body.
type Problem struct {
Type string `json:"type"`
Title string `json:"title"`
Status int `json:"status"`
Detail string `json:"detail,omitempty"`
Instance string `json:"instance,omitempty"`
}
type TaskInstanceRepository
TaskInstanceRepository reads task instances, clears them for re-run, and sets their state directly (the UI’s mark-success/failed actions).
type TaskInstanceRepository interface {
ListTaskInstances(ctx context.Context, tenant, dagID, runID string, limit, offset int) ([]domain.TaskInstance, int, error)
// ListTaskInstanceAttempts returns every attempt for (run, task), oldest
// first — the current row UNIONed with the archived history. The UI's
// /tries endpoint needs all attempts to render its navigable tabs.
ListTaskInstanceAttempts(ctx context.Context, tenant, dagID, runID, taskID string) ([]domain.TaskInstance, error)
ClearTaskInstances(ctx context.Context, tenant, dagID, runID string, taskIDs []string, onlyFailed, resetDagRun bool) (int, error)
SetTaskInstanceState(ctx context.Context, tenant, dagID, runID, taskID, state string) error
}
type TaskSummaryReader
TaskSummaryReader fetches task instances across a set of runs of a DAG, the source for the grid’s per-cell state summaries.
type TaskSummaryReader interface {
TaskInstancesForRuns(ctx context.Context, tenant, dagID string, runIDs []string) ([]domain.TaskInstance, error)
}
type TokenRenewer
TokenRenewer re-mints a still-valid user bearer with a fresh short TTL, bounded by max_lifetime since first login. *auth.JWTAuthenticator implements it via RenewUserToken; the handler depends on this narrow interface so the renew route can be tested without a real signing key.
type TokenRenewer interface {
RenewUserToken(token string, ttl, maxLifetime time.Duration) (renewed string, ok bool, err error)
}
type UIServer
UIServer serves the embedded single-page app: static assets and an index.html shell that the SPA’s client-side router falls back to. It is satisfied by internal/ui.Server. When nil, the server runs API-only and unknown paths return 404 instead of the SPA shell.
type UIServer interface {
StaticHandler() http.Handler
Index(w http.ResponseWriter, basePath string)
}
type UserAuditWriter
UserAuditWriter records account-creation events for the Audit Log. It is a separate, narrow interface (not the task-shaped AuditWriter) so account management writes a “user” resource entry with the acting admin as owner. The granted roles are passed as a single joined string so the record captures the full set.
type UserAuditWriter interface {
RecordUserCreatedAudit(ctx context.Context, tenant, actorUserID, createdUserID, email, roles string) error
}
type UserStore
UserStore creates control-plane accounts for the admin create-user API. The store hashes the plaintext password (reusing the bootstrap admin’s bcrypt scheme) and returns the created user without any secret. A duplicate email must surface as domain.ErrConflict and an unknown role as domain.ErrValidation.
type UserStore interface {
CreateUser(ctx context.Context, tenant, email, password string, roles []string) (domain.User, error)
ListUsers(ctx context.Context, tenant string, limit, offset int) ([]domain.User, int, error)
}
type VariableStore
VariableStore reads and writes Airflow-style Variables for the Admin UI.
type VariableStore interface {
ListVariables(ctx context.Context, tenant string, limit, offset int) ([]domain.Variable, int, error)
GetVariable(ctx context.Context, tenant, key string) (domain.Variable, error)
SetVariable(ctx context.Context, tenant string, v domain.Variable) error
DeleteVariable(ctx context.Context, tenant, key string) error
}
type WorkspaceFS
WorkspaceFS is the workspace-confined filesystem backing the Lite web editor (ADR 0025). Every path is relative to the workspace root and confined to it.
type WorkspaceFS interface {
Tree() ([]workspace.Entry, error)
Read(rel string) ([]byte, error)
Write(rel string, data []byte) error
Create(rel string, dir bool) error
Move(from, to string) error
Delete(rel string) error
}
type XComReader
XComReader reads stored XCom values and lists a task instance’s XCom keys for the read API.
type XComReader interface {
GetXCom(ctx context.Context, tenant, dagID, runID, taskID, key string) (xcom.Entry, error)
ListXComEntries(ctx context.Context, tenant, dagID, runID, taskID string) ([]domain.XComEntryMeta, error)
}
Generated by gomarkdoc