diff --git a/.codex b/.codex
new file mode 100644
index 000000000..e69de29bb
diff --git a/.env.local-api.example b/.env.local-api.example
new file mode 100644
index 000000000..6574e17ee
--- /dev/null
+++ b/.env.local-api.example
@@ -0,0 +1,30 @@
+APP_ENV=debug
+API_EXPORT_PORT=8029
+ROOT_API_BEARER_TOKEN=replace-with-local-dev-token
+
+DATABASE_HOST=127.0.0.1
+DATABASE_EXPORT_PORT=15432
+DATABASE_USER=replace-with-db-user
+DATABASE_PASSWORD=replace-with-db-password
+DATABASE_NAME=replace-with-db-name
+
+REDIS_HOST=127.0.0.1
+REDIS_EXPORT_PORT=16379
+REDIS_PASSWORD=replace-with-redis-password
+
+RABBITMQ_HOST=127.0.0.1
+RABBITMQ_EXPORT_PORT=15672
+RABBITMQ_USER=replace-with-rabbitmq-user
+RABBITMQ_PASSWORD=replace-with-rabbitmq-password
+RABBITMQ_VHOST=/
+RABBITMQ_VHOST_ENCODED=%2F
+
+S3_ENDPOINT=http://127.0.0.1:19000
+S3_INTERNAL_ENDPOINT=http://127.0.0.1:19000
+S3_REGION=auto
+S3_ACCESS_KEY=replace-with-s3-access-key
+S3_SECRET_KEY=replace-with-s3-secret-key
+S3_BUCKET=replace-with-s3-bucket
+
+CORE_BASE_URL=http://127.0.0.1:8019
+OTEL_EXPORTER_OTLP_ENDPOINT=
diff --git a/.gitignore b/.gitignore
index 6a7a10d44..ae49c2857 100644
--- a/.gitignore
+++ b/.gitignore
@@ -1,6 +1,9 @@
plans/
.claude
.cursorrules
+plans/
+.env.local-api
+src/server/.env.local-api
.DS_Store
.agents
skills-lock.json
diff --git a/docs/content/docs/(guides)/engineering/editing.mdx b/docs/content/docs/(guides)/engineering/editing.mdx
index 3efc2f3e1..7cfe688d9 100644
--- a/docs/content/docs/(guides)/engineering/editing.mdx
+++ b/docs/content/docs/(guides)/engineering/editing.mdx
@@ -9,32 +9,25 @@ Apply edit strategies when retrieving messages to manage context window size. Th
The `get_messages` response includes `this_time_tokens` - the total token count of returned messages. Use this to:
- Check current context window size
-- Decide when to apply edit strategies
+- Apply edit strategies only when needed
- Determine when to [reset the prompt cache](/engineering/cache)
```python title="Python"
-result = client.sessions.get_messages(session_id="session-uuid")
+result = client.sessions.get_messages(
+ session_id="session-uuid",
+ edit_strategies=[{"type": "token_limit", "params": {"limit_tokens": 30000}}],
+ editing_trigger={"token_gte": 50000},
+)
print(f"Current tokens: {result.this_time_tokens}")
-
-if result.this_time_tokens > 50000:
- # Apply strategies to reduce context
- result = client.sessions.get_messages(
- session_id="session-uuid",
- edit_strategies=[{"type": "token_limit", "params": {"limit_tokens": 30000}}]
- )
```
```typescript title="TypeScript"
-let result = await client.sessions.getMessages("session-uuid");
+const result = await client.sessions.getMessages("session-uuid", {
+ editStrategies: [{ type: "token_limit", params: { limit_tokens: 30000 } }],
+ editingTrigger: { token_gte: 50000 },
+});
console.log(`Current tokens: ${result.thisTimeTokens}`);
-
-if (result.thisTimeTokens > 50000) {
- // Apply strategies to reduce context
- result = await client.sessions.getMessages("session-uuid", {
- editStrategies: [{ type: "token_limit", params: { limit_tokens: 30000 } }],
- });
-}
```
diff --git a/src/client/acontext-py/src/acontext/_utils.py b/src/client/acontext-py/src/acontext/_utils.py
index be060649a..7d9ed6a9d 100644
--- a/src/client/acontext-py/src/acontext/_utils.py
+++ b/src/client/acontext-py/src/acontext/_utils.py
@@ -1,6 +1,6 @@
"""Utility functions for the acontext Python client."""
-from typing import Any, Iterable
+from typing import Any, Iterable, Mapping
def bool_to_str(value: bool) -> str:
@@ -58,3 +58,24 @@ def validate_edit_strategies(edit_strategies: Iterable[dict[str, Any]]) -> None:
raise ValueError("gt_token must be an integer >= 1")
if gt_token < 1:
raise ValueError("gt_token must be >= 1")
+
+
+def validate_editing_trigger(editing_trigger: Mapping[str, Any]) -> None:
+ """Validate editing trigger before sending to the API."""
+ if len(editing_trigger) == 0:
+ raise ValueError("editing_trigger must include at least one supported field")
+
+ # Keep the SDK strict so unsupported trigger names fail locally with a
+ # clearer error instead of making a round trip to the API first.
+ allowed_keys = {"token_gte"}
+ unknown_keys = set(editing_trigger.keys()) - allowed_keys
+ if unknown_keys:
+ unknown = ", ".join(sorted(unknown_keys))
+ raise ValueError(f"unsupported editing_trigger field(s): {unknown}")
+
+ if "token_gte" in editing_trigger:
+ token_gte = editing_trigger["token_gte"]
+ if isinstance(token_gte, bool) or not isinstance(token_gte, int):
+ raise ValueError("token_gte must be an integer > 0")
+ if token_gte <= 0:
+ raise ValueError("token_gte must be > 0")
diff --git a/src/client/acontext-py/src/acontext/resources/async_sessions.py b/src/client/acontext-py/src/acontext/resources/async_sessions.py
index cfd6a66bf..610f51648 100644
--- a/src/client/acontext-py/src/acontext/resources/async_sessions.py
+++ b/src/client/acontext-py/src/acontext/resources/async_sessions.py
@@ -5,12 +5,13 @@
from dataclasses import asdict
from typing import Any, BinaryIO, Literal, Optional, List
-from .._utils import build_params, validate_edit_strategies
+from .._utils import build_params, validate_edit_strategies, validate_editing_trigger
from ..client_types import AsyncRequesterProtocol
from ..messages import AcontextMessage
from ..types.common import FlagResponse
from ..types.session import (
EditStrategy,
+ EditingTrigger,
CopySessionResult,
GetMessagesOutput,
GetTasksOutput,
@@ -374,6 +375,8 @@ async def get_messages(
format: Literal["acontext", "openai", "anthropic", "gemini"] = "openai",
time_desc: bool | None = None,
edit_strategies: Optional[List[EditStrategy]] = None,
+ # editing_trigger triggers edit_strategies (v0 supports {"token_gte": int}).
+ editing_trigger: EditingTrigger | BaseModel | None = None,
pin_editing_strategies_at_message: str | None = None,
) -> GetMessagesOutput:
"""Get messages for a session.
@@ -394,6 +397,7 @@ async def get_messages(
- Middle out: [{"type": "middle_out", "params": {"token_reduce_to": 5000}}]
- Token limit: [{"type": "token_limit", "params": {"limit_tokens": 20000}}]
Defaults to None.
+ editing_trigger: Trigger config for edit_strategies, e.g. {"token_gte": 30000}. Defaults to None.
pin_editing_strategies_at_message: Message ID to pin editing strategies at.
When provided, strategies are only applied to messages up to and including
this message ID, keeping subsequent messages unchanged. This helps maintain
@@ -419,6 +423,13 @@ async def get_messages(
if edit_strategies is not None:
validate_edit_strategies(edit_strategies)
params["edit_strategies"] = json.dumps(edit_strategies)
+ if editing_trigger is not None:
+ # Keep async behavior aligned with the sync client: normalize model
+ # inputs first, then validate and serialize the exact API payload.
+ if isinstance(editing_trigger, BaseModel):
+ editing_trigger = editing_trigger.model_dump()
+ validate_editing_trigger(editing_trigger)
+ params["editing_trigger"] = json.dumps(editing_trigger)
if pin_editing_strategies_at_message is not None:
params["pin_editing_strategies_at_message"] = (
pin_editing_strategies_at_message
diff --git a/src/client/acontext-py/src/acontext/resources/sessions.py b/src/client/acontext-py/src/acontext/resources/sessions.py
index 083504a48..f61ff2d01 100644
--- a/src/client/acontext-py/src/acontext/resources/sessions.py
+++ b/src/client/acontext-py/src/acontext/resources/sessions.py
@@ -5,12 +5,13 @@
from dataclasses import asdict
from typing import Any, BinaryIO, Literal, Optional, List
-from .._utils import build_params, validate_edit_strategies
+from .._utils import build_params, validate_edit_strategies, validate_editing_trigger
from ..client_types import RequesterProtocol
from ..messages import AcontextMessage
from ..types.common import FlagResponse
from ..types.session import (
EditStrategy,
+ EditingTrigger,
CopySessionResult,
GetMessagesOutput,
GetTasksOutput,
@@ -374,6 +375,8 @@ def get_messages(
format: Literal["acontext", "openai", "anthropic", "gemini"] = "openai",
time_desc: bool | None = None,
edit_strategies: Optional[List[EditStrategy]] = None,
+ # editing_trigger triggers edit_strategies (v0 supports {"token_gte": int}).
+ editing_trigger: EditingTrigger | BaseModel | None = None,
pin_editing_strategies_at_message: str | None = None,
) -> GetMessagesOutput:
"""Get messages for a session.
@@ -394,6 +397,7 @@ def get_messages(
- Middle out: [{"type": "middle_out", "params": {"token_reduce_to": 5000}}]
- Token limit: [{"type": "token_limit", "params": {"limit_tokens": 20000}}]
Defaults to None.
+ editing_trigger: Trigger config for edit_strategies, e.g. {"token_gte": 30000}. Defaults to None.
pin_editing_strategies_at_message: Message ID to pin editing strategies at.
When provided, strategies are only applied to messages up to and including
this message ID, keeping subsequent messages unchanged. This helps maintain
@@ -419,6 +423,13 @@ def get_messages(
if edit_strategies is not None:
validate_edit_strategies(edit_strategies)
params["edit_strategies"] = json.dumps(edit_strategies)
+ if editing_trigger is not None:
+ # Accept either a plain dict or a caller-provided Pydantic model so
+ # the SDK surface matches the existing flexibility of edit_strategies.
+ if isinstance(editing_trigger, BaseModel):
+ editing_trigger = editing_trigger.model_dump()
+ validate_editing_trigger(editing_trigger)
+ params["editing_trigger"] = json.dumps(editing_trigger)
if pin_editing_strategies_at_message is not None:
params["pin_editing_strategies_at_message"] = (
pin_editing_strategies_at_message
diff --git a/src/client/acontext-py/src/acontext/types/__init__.py b/src/client/acontext-py/src/acontext/types/__init__.py
index 28370a6bd..ba3797b65 100644
--- a/src/client/acontext-py/src/acontext/types/__init__.py
+++ b/src/client/acontext-py/src/acontext/types/__init__.py
@@ -12,6 +12,7 @@
)
from .session import (
Asset,
+ EditingTrigger,
GetMessagesOutput,
GetTasksOutput,
ListSessionsOutput,
@@ -67,6 +68,7 @@
"UpdateArtifactResp",
# Session types
"Asset",
+ "EditingTrigger",
"GetMessagesOutput",
"GetTasksOutput",
"ListSessionsOutput",
diff --git a/src/client/acontext-py/src/acontext/types/session.py b/src/client/acontext-py/src/acontext/types/session.py
index 8f07b787c..b7277dacd 100644
--- a/src/client/acontext-py/src/acontext/types/session.py
+++ b/src/client/acontext-py/src/acontext/types/session.py
@@ -126,6 +126,17 @@ class MiddleOutStrategy(TypedDict):
]
+class EditingTrigger(TypedDict, total=False):
+ """Trigger config for applying edit strategies.
+
+ Attributes:
+ token_gte: Apply edit strategies only when the current token count is
+ greater than or equal to this value.
+ """
+
+ token_gte: NotRequired[int]
+
+
class Asset(BaseModel):
"""Asset model representing a file asset."""
diff --git a/src/client/acontext-ts/src/resources/sessions.ts b/src/client/acontext-ts/src/resources/sessions.ts
index 2ff2aacdd..4825dc319 100644
--- a/src/client/acontext-ts/src/resources/sessions.ts
+++ b/src/client/acontext-ts/src/resources/sessions.ts
@@ -9,6 +9,8 @@ import { buildParams, validateUUID } from '../utils';
import {
EditStrategy,
EditStrategySchema,
+ EditingTrigger,
+ EditingTriggerSchema,
CopySessionResult,
CopySessionResultSchema,
FlagResponse,
@@ -340,6 +342,7 @@ export class SessionsAPI {
* @param options.format - The format of the messages ('acontext', 'openai', 'anthropic', or 'gemini').
* @param options.timeDesc - Order by created_at descending if true, ascending if false.
* @param options.editStrategies - Optional list of edit strategies to apply before format conversion.
+ * @param options.editingTrigger - Optional trigger config for editStrategies (v0 supports { token_gte: number }).
* Examples:
* - Remove tool results: [{ type: 'remove_tool_result', params: { keep_recent_n_tool_results: 3 } }]
* - Remove large tool results: [{ type: 'remove_tool_result', params: { gt_token: 100 } }]
@@ -364,6 +367,7 @@ export class SessionsAPI {
format?: 'acontext' | 'openai' | 'anthropic' | 'gemini';
timeDesc?: boolean | null;
editStrategies?: Array | null;
+ editingTrigger?: EditingTrigger | null;
pinEditingStrategiesAtMessage?: string | null;
}
): Promise {
@@ -387,6 +391,12 @@ export class SessionsAPI {
EditStrategySchema.array().parse(options.editStrategies);
params.edit_strategies = JSON.stringify(options.editStrategies);
}
+ if (options?.editingTrigger !== undefined && options?.editingTrigger !== null) {
+ // Validate before serializing so unsupported trigger shapes fail at the
+ // SDK boundary rather than after an API request.
+ EditingTriggerSchema.parse(options.editingTrigger);
+ params.editing_trigger = JSON.stringify(options.editingTrigger);
+ }
if (options?.pinEditingStrategiesAtMessage !== undefined && options?.pinEditingStrategiesAtMessage !== null) {
params.pin_editing_strategies_at_message = options.pinEditingStrategiesAtMessage;
}
diff --git a/src/client/acontext-ts/src/types/session.ts b/src/client/acontext-ts/src/types/session.ts
index 02951f1d5..42ac8e925 100644
--- a/src/client/acontext-ts/src/types/session.ts
+++ b/src/client/acontext-ts/src/types/session.ts
@@ -328,3 +328,17 @@ export const EditStrategySchema = z.union([
]);
export type EditStrategy = z.infer;
+
+/**
+ * Trigger config for applying edit strategies.
+ * v0 supports only token_gte.
+ */
+export const EditingTriggerSchema = z.object({
+ token_gte: z.number().int().positive().optional(),
+}).strict().refine((value) => Object.keys(value).length > 0, {
+ // Mirror the API's "at least one supported trigger" rule so empty objects
+ // are rejected consistently across clients and server.
+ message: 'editingTrigger must include at least one supported field',
+});
+
+export type EditingTrigger = z.infer;
diff --git a/src/server/.env.local-api.example b/src/server/.env.local-api.example
new file mode 100644
index 000000000..6574e17ee
--- /dev/null
+++ b/src/server/.env.local-api.example
@@ -0,0 +1,30 @@
+APP_ENV=debug
+API_EXPORT_PORT=8029
+ROOT_API_BEARER_TOKEN=replace-with-local-dev-token
+
+DATABASE_HOST=127.0.0.1
+DATABASE_EXPORT_PORT=15432
+DATABASE_USER=replace-with-db-user
+DATABASE_PASSWORD=replace-with-db-password
+DATABASE_NAME=replace-with-db-name
+
+REDIS_HOST=127.0.0.1
+REDIS_EXPORT_PORT=16379
+REDIS_PASSWORD=replace-with-redis-password
+
+RABBITMQ_HOST=127.0.0.1
+RABBITMQ_EXPORT_PORT=15672
+RABBITMQ_USER=replace-with-rabbitmq-user
+RABBITMQ_PASSWORD=replace-with-rabbitmq-password
+RABBITMQ_VHOST=/
+RABBITMQ_VHOST_ENCODED=%2F
+
+S3_ENDPOINT=http://127.0.0.1:19000
+S3_INTERNAL_ENDPOINT=http://127.0.0.1:19000
+S3_REGION=auto
+S3_ACCESS_KEY=replace-with-s3-access-key
+S3_SECRET_KEY=replace-with-s3-secret-key
+S3_BUCKET=replace-with-s3-bucket
+
+CORE_BASE_URL=http://127.0.0.1:8019
+OTEL_EXPORTER_OTLP_ENDPOINT=
diff --git a/src/server/api/go/internal/modules/handler/session.go b/src/server/api/go/internal/modules/handler/session.go
index f8237ecd7..8263a08ce 100644
--- a/src/server/api/go/internal/modules/handler/session.go
+++ b/src/server/api/go/internal/modules/handler/session.go
@@ -19,6 +19,7 @@ import (
"github.com/memodb-io/Acontext/internal/modules/serializer"
"github.com/memodb-io/Acontext/internal/modules/service"
"github.com/memodb-io/Acontext/internal/pkg/converter"
+ "github.com/memodb-io/Acontext/internal/pkg/editingtrigger"
"github.com/memodb-io/Acontext/internal/pkg/editor"
"github.com/memodb-io/Acontext/internal/pkg/normalizer"
"github.com/memodb-io/Acontext/internal/pkg/tokenizer"
@@ -495,6 +496,7 @@ type GetMessagesReq struct {
Format string `form:"format,default=openai" json:"format" binding:"omitempty,oneof=acontext openai anthropic gemini" example:"openai" enums:"acontext,openai,anthropic,gemini"`
TimeDesc bool `form:"time_desc,default=false" json:"time_desc" example:"false"`
EditStrategies string `form:"edit_strategies" json:"edit_strategies" example:"[{\"type\":\"remove_tool_result\",\"params\":{\"keep_recent_n_tool_results\":3}}]"`
+ EditingTrigger string `form:"editing_trigger" json:"editing_trigger" example:"{\"token_gte\":30000}"`
PinEditingStrategiesAtMessage string `form:"pin_editing_strategies_at_message" json:"pin_editing_strategies_at_message" example:""`
}
@@ -513,6 +515,7 @@ type GetMessagesReq struct {
// @Param format query string false "Format to convert messages to: acontext (original), openai (default), anthropic, gemini." enums(acontext,openai,anthropic,gemini)
// @Param time_desc query boolean false "Order by created_at descending if true, ascending if false (default false)" example(false)
// @Param edit_strategies query string false "JSON array of edit strategies to apply before format conversion" example([{"type":"remove_tool_result","params":{"keep_recent_n_tool_results":3}}])
+// @Param editing_trigger query string false "JSON object trigger for edit_strategies. v0 supports only {\"token_gte\": } (OR semantics when more triggers are added)." example({"token_gte":30000})
// @Param pin_editing_strategies_at_message query string false "Message ID to pin editing strategies at. When provided, strategies are only applied to messages up to and including this message ID, keeping subsequent messages unchanged. This helps maintain prompt cache stability by preserving a stable prefix. The response will include edit_at_message_id indicating where strategies were applied." example()
// @Security BearerAuth
// @Success 200 {object} serializer.Response{data=converter.GetMessagesOutput}
@@ -552,6 +555,40 @@ func (h *SessionHandler) GetMessages(c *gin.Context) {
}
}
+ // Parse editing_trigger if provided (v0 supports only token_gte).
+ var editingTrigger *service.EditingTrigger
+ if req.EditingTrigger != "" {
+ // editing_trigger is only meaningful when there is something to gate.
+ // Rejecting it here keeps the API contract explicit instead of silently
+ // accepting a no-op parameter.
+ if req.EditStrategies == "" {
+ c.JSON(http.StatusBadRequest, serializer.ParamErr("editing_trigger requires edit_strategies", errors.New("missing edit_strategies")))
+ return
+ }
+
+ var trig service.EditingTrigger
+ if err := json.Unmarshal([]byte(req.EditingTrigger), &trig); err != nil {
+ // Surface "unknown field" separately so callers can distinguish
+ // unsupported trigger names from malformed JSON payloads.
+ var unsupportedErr editingtrigger.UnsupportedTriggerError
+ if errors.As(err, &unsupportedErr) {
+ c.JSON(http.StatusBadRequest, serializer.ParamErr("invalid editing_trigger", err))
+ return
+ }
+ c.JSON(http.StatusBadRequest, serializer.ParamErr("invalid editing_trigger JSON", err))
+ return
+ }
+ if err := trig.Validate(); err != nil {
+ if errors.Is(err, editingtrigger.ErrTokenGteMustBeGreater) {
+ c.JSON(http.StatusBadRequest, serializer.ParamErr("invalid editing_trigger.token_gte", err))
+ return
+ }
+ c.JSON(http.StatusBadRequest, serializer.ParamErr("invalid editing_trigger", err))
+ return
+ }
+ editingTrigger = &trig
+ }
+
out, err := h.svc.GetMessages(c.Request.Context(), service.GetMessagesInput{
ProjectID: project.ID,
SessionID: sessionID,
@@ -562,10 +599,18 @@ func (h *SessionHandler) GetMessages(c *gin.Context) {
AssetExpire: time.Hour * 24,
TimeDesc: req.TimeDesc,
EditStrategies: editStrategies,
+ EditingTrigger: editingTrigger,
PinEditingStrategiesAtMessage: req.PinEditingStrategiesAtMessage,
UserKEK: middleware.GetUserKEKIfEncrypted(c),
})
if err != nil {
+ // Token counting now happens inside the service because trigger evaluation
+ // and final response tokens share that result. Promote those failures to
+ // 500 so they are treated as server-side computation errors.
+ if errors.Is(err, service.ErrGetMessagesTokenCount) {
+ c.JSON(http.StatusInternalServerError, serializer.DBErr("failed to count tokens", err))
+ return
+ }
c.JSON(http.StatusBadRequest, serializer.DBErr("", err))
return
}
@@ -582,13 +627,6 @@ func (h *SessionHandler) GetMessages(c *gin.Context) {
return
}
- // Calculate token count for the returned messages
- thisTimeTokens, err := tokenizer.CountMessagePartsTokens(c.Request.Context(), out.Items)
- if err != nil {
- c.JSON(http.StatusInternalServerError, serializer.DBErr("failed to count tokens", err))
- return
- }
-
convertedOut, err := converter.GetConvertedMessagesOutput(
out.Items,
format,
@@ -596,7 +634,7 @@ func (h *SessionHandler) GetMessages(c *gin.Context) {
out.Events,
out.NextCursor,
out.HasMore,
- thisTimeTokens,
+ out.ThisTimeTokens,
out.EditAtMessageID,
)
if err != nil {
diff --git a/src/server/api/go/internal/modules/service/session.go b/src/server/api/go/internal/modules/service/session.go
index d1fa0a33f..75dd77261 100644
--- a/src/server/api/go/internal/modules/service/session.go
+++ b/src/server/api/go/internal/modules/service/session.go
@@ -2,10 +2,13 @@ package service
import (
"context"
+ "crypto/sha256"
"encoding/base64"
"encoding/binary"
+ "encoding/hex"
"errors"
"fmt"
+ "io"
"mime/multipart"
"sort"
"time"
@@ -19,8 +22,10 @@ import (
mq "github.com/memodb-io/Acontext/internal/infra/queue"
"github.com/memodb-io/Acontext/internal/modules/model"
"github.com/memodb-io/Acontext/internal/modules/repo"
+ "github.com/memodb-io/Acontext/internal/pkg/editingtrigger"
"github.com/memodb-io/Acontext/internal/pkg/editor"
"github.com/memodb-io/Acontext/internal/pkg/paging"
+ "github.com/memodb-io/Acontext/internal/pkg/tokenizer"
"github.com/redis/go-redis/v9"
"go.uber.org/zap"
"gorm.io/datatypes"
@@ -67,6 +72,8 @@ type sessionService struct {
materialSvc MaterialService
}
+var ErrGetMessagesTokenCount = errors.New("get messages token count error")
+
const (
// Redis key prefix for message parts cache
redisKeyPrefixParts = "message:parts:"
@@ -436,19 +443,23 @@ func (s *sessionService) StoreMessage(ctx context.Context, in StoreMessageInput)
}
type GetMessagesInput struct {
- ProjectID uuid.UUID `json:"project_id"`
- SessionID uuid.UUID `json:"session_id"`
- Limit int `json:"limit"`
- Cursor string `json:"cursor"`
- WithAssetPublicURL bool `json:"with_public_url"`
- AssetExpire time.Duration `json:"asset_expire"`
- TimeDesc bool `json:"time_desc"`
- WithEvents bool `json:"with_events"`
- EditStrategies []editor.StrategyConfig `json:"edit_strategies,omitempty"`
- PinEditingStrategiesAtMessage string `json:"pin_editing_strategies_at_message,omitempty"`
- UserKEK []byte `json:"-"` // optional: for envelope encryption (decrypting parts)
+ ProjectID uuid.UUID `json:"project_id"`
+ SessionID uuid.UUID `json:"session_id"`
+ Limit int `json:"limit"`
+ Cursor string `json:"cursor"`
+ WithAssetPublicURL bool `json:"with_public_url"`
+ AssetExpire time.Duration `json:"asset_expire"`
+ TimeDesc bool `json:"time_desc"`
+ WithEvents bool `json:"with_events"`
+ EditStrategies []editor.StrategyConfig `json:"edit_strategies,omitempty"`
+ // EditingTrigger holds optional trigger config for applying edit_strategies.
+ EditingTrigger *EditingTrigger `json:"editing_trigger,omitempty"`
+ PinEditingStrategiesAtMessage string `json:"pin_editing_strategies_at_message,omitempty"`
+ UserKEK []byte `json:"-"` // optional: for envelope encryption (decrypting parts)
}
+type EditingTrigger = editingtrigger.Trigger
+
type PublicURL struct {
URL string `json:"url"`
ExpireAt time.Time `json:"expire_at"`
@@ -460,6 +471,7 @@ type GetMessagesOutput struct {
NextCursor string `json:"next_cursor,omitempty"`
HasMore bool `json:"has_more"`
PublicURLs map[string]PublicURL `json:"public_urls,omitempty"` // file_name -> url
+ ThisTimeTokens int `json:"this_time_tokens"`
EditAtMessageID string `json:"edit_at_message_id,omitempty"`
}
@@ -551,13 +563,75 @@ func (s *sessionService) GetMessages(ctx context.Context, in GetMessagesInput) (
}
// Apply edit strategies if provided (before format conversion)
+ var triggerEval *editingtrigger.Eval
+ strategiesApplied := false
if len(in.EditStrategies) > 0 {
- result, err := editor.ApplyStrategiesWithPin(out.Items, in.EditStrategies, in.PinEditingStrategiesAtMessage)
- if err != nil {
- return nil, fmt.Errorf("failed to apply edit strategies: %w", err)
+ // Preserve the previous behavior by default: strategies run whenever they
+ // are provided. editing_trigger only changes that default when present.
+ applyEditStrategies := true
+ triggerChecks := editingtrigger.BuildChecks(in.EditingTrigger)
+ triggerEvaluated := len(triggerChecks) > 0
+ if triggerEvaluated {
+ // Evaluate trigger on the same editable prefix used by pin_editing_strategies_at_message.
+ triggerMessages := out.Items
+ if in.PinEditingStrategiesAtMessage != "" {
+ pinIndex := -1
+ for i := range out.Items {
+ if out.Items[i].ID.String() == in.PinEditingStrategiesAtMessage {
+ pinIndex = i
+ break
+ }
+ }
+ if pinIndex != -1 {
+ triggerMessages = out.Items[:pinIndex+1]
+ }
+ }
+
+ // OR semantics: apply when any trigger check passes.
+ applyEditStrategies = false
+ eval := editingtrigger.NewEval(
+ in.SessionID,
+ triggerMessages,
+ func(ctx context.Context, messages []model.Message) (int, error) {
+ // Wrap tokenizer failures with a sentinel so the handler can map
+ // trigger/token computation failures to a stable HTTP 500 response.
+ tokens, err := tokenizer.CountMessagePartsTokens(ctx, messages)
+ if err != nil {
+ return 0, fmt.Errorf("%w: failed to count tokens for editing_trigger session_id=%s: %v", ErrGetMessagesTokenCount, in.SessionID, err)
+ }
+ return tokens, nil
+ },
+ )
+ triggerEval = eval
+ for _, check := range triggerChecks {
+ ok, err := check(ctx, eval)
+ if err != nil {
+ return nil, err
+ }
+ if ok {
+ applyEditStrategies = true
+ break
+ }
+ }
+
+ }
+
+ if applyEditStrategies {
+ result, err := editor.ApplyStrategiesWithPin(out.Items, in.EditStrategies, in.PinEditingStrategiesAtMessage)
+ if err != nil {
+ return nil, fmt.Errorf("failed to apply edit strategies: %w", err)
+ }
+ strategiesApplied = true
+ out.Items = result.Messages
+ out.EditAtMessageID = result.EditAtMessageID
+ } else if triggerEvaluated && in.PinEditingStrategiesAtMessage != "" {
+ // Trigger skipped editing; preserve caller-provided boundary for future requests.
+ out.EditAtMessageID = in.PinEditingStrategiesAtMessage
+ } else if out.EditAtMessageID == "" && len(out.Items) > 0 {
+ // Even when editing is skipped, return a deterministic edit boundary so
+ // clients can reuse the latest message ID in the next request.
+ out.EditAtMessageID = out.Items[len(out.Items)-1].ID.String()
}
- out.Items = result.Messages
- out.EditAtMessageID = result.EditAtMessageID
} else if len(out.Items) > 0 {
// No strategies, but still set EditAtMessageID to the last message
out.EditAtMessageID = out.Items[len(out.Items)-1].ID.String()
@@ -588,6 +662,27 @@ func (s *sessionService) GetMessages(ctx context.Context, in GetMessagesInput) (
}
}
+ usedCachedTokens := false
+ if triggerEval != nil {
+ // Reuse the trigger-time token count only when the final output is
+ // byte-for-byte equivalent for token purposes. Editing can mutate
+ // message content even if IDs stay unchanged, so guard reuse carefully.
+ if cachedTokens, ok := triggerEval.CachedTokens(); ok && !strategiesApplied && sameMessageContentSignature(triggerEval.Messages(), out.Items) {
+ out.ThisTimeTokens = cachedTokens
+ usedCachedTokens = true
+ }
+ }
+
+ if !usedCachedTokens {
+ // Fall back to counting the final returned payload so this_time_tokens
+ // always reflects what the client actually receives.
+ thisTimeTokens, err := tokenizer.CountMessagePartsTokens(ctx, out.Items)
+ if err != nil {
+ return nil, fmt.Errorf("%w: session_id=%s: %v", ErrGetMessagesTokenCount, in.SessionID, err)
+ }
+ out.ThisTimeTokens = thisTimeTokens
+ }
+
return out, nil
}
@@ -599,6 +694,67 @@ func (s *sessionService) DownloadAsset(ctx context.Context, s3Key string, userKE
return s.s3.DownloadFile(ctx, s3Key, userKEK)
}
+// sameMessageContentSignature returns true only when two message slices are
+// equivalent for token-count reuse.
+//
+// Why this is needed:
+// sameMessageOrderByID was not enough, because token count depends on content.
+// IDs can stay the same while token-relevant content changes.
+//
+// Example:
+// - triggerEval.Messages():
+// - msg-1 text: "hello"
+// - msg-2 tool-call: name="search", arguments={"q":"apple"}
+// - out.Items (same msg IDs/order):
+// - msg-1 text: "hello"
+// - msg-2 tool-call: name="search", arguments={"q":"banana"}
+//
+// ID-only compare would return true and reuse stale cached tokens.
+// Content-signature compare returns false, so token count is recomputed.
+func sameMessageContentSignature(a, b []model.Message) bool {
+ sigA, err := messageTokenSignature(a)
+ if err != nil {
+ return false
+ }
+ sigB, err := messageTokenSignature(b)
+ if err != nil {
+ return false
+ }
+ return sigA == sigB
+}
+
+func messageTokenSignature(messages []model.Message) (string, error) {
+ hasher := sha256.New()
+ if _, err := io.WriteString(hasher, fmt.Sprintf("%d|", len(messages))); err != nil {
+ return "", err
+ }
+
+ for _, msg := range messages {
+ if _, err := io.WriteString(hasher, msg.ID.String()); err != nil {
+ return "", err
+ }
+ if _, err := io.WriteString(hasher, "|"); err != nil {
+ return "", err
+ }
+
+ // Hash only the token-relevant projection of a message. That keeps the
+ // reuse check aligned with tokenizer behavior instead of unrelated fields
+ // like timestamps or DB metadata.
+ content, err := tokenizer.ExtractTextAndToolContent(msg.Parts)
+ if err != nil {
+ return "", err
+ }
+ if _, err := io.WriteString(hasher, content); err != nil {
+ return "", err
+ }
+ if _, err := io.WriteString(hasher, "\n---\n"); err != nil {
+ return "", err
+ }
+ }
+
+ return hex.EncodeToString(hasher.Sum(nil)), nil
+}
+
// cachePartsInRedis stores message parts in Redis with a fixed TTL.
// When userKEK is provided, the serialized JSON is encrypted before caching.
// Format: prefix_byte | payload
diff --git a/src/server/api/go/internal/pkg/editingtrigger/editing_trigger.go b/src/server/api/go/internal/pkg/editingtrigger/editing_trigger.go
new file mode 100644
index 000000000..dfc7b6544
--- /dev/null
+++ b/src/server/api/go/internal/pkg/editingtrigger/editing_trigger.go
@@ -0,0 +1,180 @@
+package editingtrigger
+
+import (
+ "context"
+ "encoding/json"
+ "errors"
+ "fmt"
+
+ "github.com/google/uuid"
+ "github.com/memodb-io/Acontext/internal/modules/model"
+)
+
+// Trigger defines trigger configuration for applying edit strategies.
+// v0 supports only token_gte.
+type Trigger struct {
+ // TokenGte triggers edit strategies when token count is >= this value.
+ TokenGte *int `json:"token_gte,omitempty"`
+
+ rawKeys map[string]struct{} `json:"-"`
+}
+
+var (
+ ErrNoSupportedTrigger = errors.New("at least one supported trigger is required")
+ ErrTokenGteMustBeGreater = errors.New("token_gte must be > 0")
+)
+
+type UnsupportedTriggerError struct {
+ Key string
+}
+
+func (e UnsupportedTriggerError) Error() string {
+ return fmt.Sprintf("unsupported trigger: %s", e.Key)
+}
+
+func (t *Trigger) UnmarshalJSON(data []byte) error {
+ var raw map[string]json.RawMessage
+ if err := json.Unmarshal(data, &raw); err != nil {
+ return err
+ }
+
+ t.TokenGte = nil
+ // Track which keys were explicitly present so validation can distinguish
+ // between "field omitted" and "field provided with a bad/null value".
+ t.rawKeys = make(map[string]struct{}, len(raw))
+
+ for key, value := range raw {
+ switch key {
+ case "token_gte":
+ t.rawKeys[key] = struct{}{}
+ if err := json.Unmarshal(value, &t.TokenGte); err != nil {
+ return fmt.Errorf("invalid token_gte: %w", err)
+ }
+ default:
+ return UnsupportedTriggerError{Key: key}
+ }
+ }
+
+ return nil
+}
+
+func (t Trigger) Validate() error {
+ hasAnySupportedTrigger := len(t.rawKeys) > 0 || t.TokenGte != nil
+ if !hasAnySupportedTrigger {
+ return ErrNoSupportedTrigger
+ }
+
+ // {"token_gte": null} unmarshals to a nil pointer, so use rawKeys to keep
+ // rejecting it instead of treating it as "not configured".
+ _, tokenGteProvided := t.rawKeys["token_gte"]
+ if tokenGteProvided && t.TokenGte == nil {
+ return ErrTokenGteMustBeGreater
+ }
+ if t.TokenGte != nil && *t.TokenGte <= 0 {
+ return ErrTokenGteMustBeGreater
+ }
+
+ return nil
+}
+
+// TokenCounter computes token count for a message slice.
+type TokenCounter func(ctx context.Context, messages []model.Message) (int, error)
+
+// Eval evaluates trigger checks and memoizes token count.
+type Eval struct {
+ sessionID uuid.UUID
+ messages []model.Message
+ counter TokenCounter
+
+ tokenCount *int
+}
+
+func NewEval(sessionID uuid.UUID, messages []model.Message, counter TokenCounter) *Eval {
+ return &Eval{
+ sessionID: sessionID,
+ messages: messages,
+ counter: counter,
+ }
+}
+
+func (e *Eval) Tokens(ctx context.Context) (int, error) {
+ if e.tokenCount != nil {
+ // Trigger checks can ask for tokens more than once; memoize so multiple
+ // checks still pay the tokenizer cost only once per request.
+ return *e.tokenCount, nil
+ }
+
+ tokens, err := e.counter(ctx, e.messages)
+ if err != nil {
+ return 0, err
+ }
+
+ e.tokenCount = &tokens
+ return tokens, nil
+}
+
+func (e *Eval) Messages() []model.Message {
+ return e.messages
+}
+
+func (e *Eval) CachedTokens() (int, bool) {
+ if e.tokenCount == nil {
+ return 0, false
+ }
+ return *e.tokenCount, true
+}
+
+type Check func(ctx context.Context, eval *Eval) (bool, error)
+
+type namedCheck struct {
+ name string
+ build func(trigger *Trigger) []Check
+}
+
+type registry struct {
+ checks []namedCheck
+}
+
+var triggerRegistry = registry{
+ checks: []namedCheck{
+ {
+ name: "token_gte",
+ build: tokenGteChecks,
+ },
+ },
+}
+
+func tokenGteChecks(trigger *Trigger) []Check {
+ if trigger.TokenGte == nil || *trigger.TokenGte <= 0 {
+ return nil
+ }
+
+ threshold := *trigger.TokenGte
+ return []Check{
+ func(ctx context.Context, eval *Eval) (bool, error) {
+ tokens, err := eval.Tokens(ctx)
+ if err != nil {
+ return false, err
+ }
+ // token_gte is intentionally inclusive so callers can pin a hard
+ // threshold without off-by-one ambiguity.
+ return tokens >= threshold, nil
+ },
+ }
+}
+
+func BuildChecks(trigger *Trigger) []Check {
+ if trigger == nil {
+ return nil
+ }
+
+ // Build a flat list of checks so the service can evaluate them with OR
+ // semantics while keeping trigger registration centralized here.
+ checks := make([]Check, 0, len(triggerRegistry.checks))
+ for _, entry := range triggerRegistry.checks {
+ _ = entry.name
+ checks = append(checks, entry.build(trigger)...)
+ }
+
+ return checks
+}