Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
953 changes: 534 additions & 419 deletions clientlibrary/service/service.pb.go

Large diffs are not rendered by default.

9 changes: 9 additions & 0 deletions core/anytype/bootstrap.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import (
"github.com/anyproto/any-sync/commonfile/fileservice"
"github.com/anyproto/any-sync/commonspace"
"github.com/anyproto/any-sync/commonspace/acl/aclclient"
anysyncpubsub "github.com/anyproto/any-sync/commonspace/pubsub"
anysyncinboxclient "github.com/anyproto/any-sync/coordinator/inboxclient"

"github.com/anyproto/any-sync/coordinator/nodeconfsource"
Expand Down Expand Up @@ -103,6 +104,7 @@ import (
"github.com/anyproto/anytype-heart/core/payments/emailcollector"
"github.com/anyproto/anytype-heart/core/peerstatus"
"github.com/anyproto/anytype-heart/core/publish"
"github.com/anyproto/anytype-heart/core/pubsub"
"github.com/anyproto/anytype-heart/core/pushnotification"
"github.com/anyproto/anytype-heart/core/pushnotification/pushclient"
"github.com/anyproto/anytype-heart/core/relationutils/formatfetcher"
Expand Down Expand Up @@ -240,6 +242,10 @@ func Bootstrap(a *app.App, components ...app.Component) {
a.Register(c)
}

// the pubsub engine takes its client-side deps (crypto, membership, peers)
// from the heart-side component, resolved lazily after Init
pubsubService := pubsub.New()

a.
// profiler is registered early so its Init wires the event sender
// before any storage component runs its own Init — storage corruption
Expand Down Expand Up @@ -299,6 +305,9 @@ func Bootstrap(a *app.App, components ...app.Component) {
Register(transportpenalty.New()).
Register(localdiscovery.New()).
Register(peermanager.New()).
Register(anysyncpubsub.New(pubsubService.EngineDeps())).
// Cancel accepted publishes before the engine waits for its workers.
Register(pubsubService).
Register(typeprovider.New()).
Register(fileuploader.New()).
Register(rpcstore.New()).
Expand Down
1 change: 1 addition & 0 deletions core/api/core/core.go
Original file line number Diff line number Diff line change
Expand Up @@ -265,6 +265,7 @@ type ClientCommands interface {
ChatReadMessages(context.Context, *pb.RpcChatReadMessagesRequest) *pb.RpcChatReadMessagesResponse
ChatReadReactions(context.Context, *pb.RpcChatReadReactionsRequest) *pb.RpcChatReadReactionsResponse
ChatSearch(context.Context, *pb.RpcChatSearchRequest) *pb.RpcChatSearchResponse
PubsubPublish(context.Context, *pb.RpcPubsubPublishRequest) *pb.RpcPubsubPublishResponse
}

// WidgetScope names which sidebar root a widget lives in. A space has two:
Expand Down
49 changes: 49 additions & 0 deletions core/api/core/mock_apicore/mock_ClientCommands.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

140 changes: 139 additions & 1 deletion core/api/docs/v2/openapi.json

Large diffs are not rendered by default.

121 changes: 121 additions & 0 deletions core/api/docs/v2/openapi.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -575,6 +575,11 @@ components:
unread_reaction_order:
type: string
type: object
ChatStatusResult:
properties:
dry_run:
type: boolean
type: object
CreateChatRequest:
properties:
name:
Expand Down Expand Up @@ -1638,6 +1643,35 @@ info:
A read does not fail because it encounters content its representation
cannot express. Such content is reported in `warnings` beside the result.

## Publish chat status

`POST /v2/spaces/{space_id}/chats/{chat_id}/status` sends ephemeral activity
for a chat or discussion. For a typing indicator, send an empty body or `{}`.
For detailed agent activity, send:

```json
{"text":"Searching documentation","data":{"tool_call":"web_search"}}
```

`text` is optional display text. Empty or omitted text stays omitted on the
wire, allowing clients to show their localized typing label. `data` is optional
and accepts any JSON value, including arrays, scalars, and null. The encoded
JSON payload, including field names and escaping, must fit in 65,508 bytes.

Success is `200 {}` and means the update was accepted for publication. The
Space's pubsub topic is `status/<chat_id>`; subscribers receive
`Event.Pubsub.Message` with the verified sender identity. Status is neither
stored as a chat message nor replayed by the chat-message stream.

This route requires write access. `?dry_run=true` validates without publishing;
`Idempotency-Key` retries replay the response without publishing again. Use
a fresh key for each refresh, or omit it.

Recommended client convention: refresh every two seconds while active and
expire the status ten seconds after its last receipt. Stop refreshing when
finished. Expiry is the receiving client's responsibility; empty text still
means activity and does not clear it.

## Stream chat messages

`GET /v2/spaces/{space_id}/chats/{chat_id}/messages/stream` opens a
Expand Down Expand Up @@ -2835,6 +2869,93 @@ paths:
summary: Mark chat activity as read
tags:
- Chat
/v2/spaces/{space_id}/chats/{chat_id}/status:
post:
description: Publishes encrypted ephemeral JSON on status/<chat_id>. Empty text
is omitted so clients can show localized typing. data accepts any JSON value.
The encoded payload is limited to 65,508 bytes. Status is not stored or replayed;
receivers expire it. Success acknowledges acceptance, not delivery.
operationId: publish_chat_status
parameters:
- description: Space id
in: path
name: space_id
required: true
schema:
type: string
- description: Chat or discussion id
in: path
name: chat_id
required: true
schema:
type: string
- description: Validate without publishing
in: query
name: dry_run
schema:
type: boolean
- description: Unique key for this update; a reused key replays the response
without publishing again
in: header
name: Idempotency-Key
schema:
type: string
requestBody:
content:
application/json:
schema:
additionalProperties: false
description: ephemeral activity; the encoded JSON payload must fit in 65508 bytes, including field names and escaping
properties:
data:
description: 'optional arbitrary JSON value: object, array, string, number, boolean or null'
text:
description: optional display text; empty or omitted is left out of the pubsub payload so receivers can show a localized typing label; JSON escaping and data consume the same byte budget
maxLength: 65497
type: string
type: object
description: Optional text and arbitrary JSON data; clients localize the default
typing label
responses:
"200":
content:
application/json:
schema:
$ref: '#/components/schemas/ChatStatusResult'
description: Status accepted, or validated on a dry run
"400":
content:
application/json:
schema:
$ref: '#/components/schemas/Error'
description: Invalid status
"401":
$ref: '#/components/responses/Unauthorized'
"403":
$ref: '#/components/responses/Forbidden'
"404":
content:
application/json:
schema:
$ref: '#/components/schemas/Error'
description: Space or chat not found
"409":
$ref: '#/components/responses/Conflict'
"413":
$ref: '#/components/responses/RequestTooLarge'
"429":
$ref: '#/components/responses/RateLimited'
"500":
content:
application/json:
schema:
$ref: '#/components/schemas/Error'
description: Publication failed
security:
- bearerauth: []
summary: Publish chat status
tags:
- Chat
/v2/spaces/{space_id}/collections:
post:
description: Item ids are checked against the space; an id that does not resolve
Expand Down
2 changes: 2 additions & 0 deletions core/api/server/v2_router_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -271,6 +271,7 @@ func TestV2Routes(t *testing.T) {
{"DELETE", "/v2/spaces/space1/chats/chat1/messages/msg1"},
{"POST", "/v2/spaces/space1/chats/chat1/messages/msg1/reactions"},
{"POST", "/v2/spaces/space1/chats/chat1/read"},
{"POST", "/v2/spaces/space1/chats/chat1/status"},
// the sidebar mutations: a retried create duplicates a widget the
// desktop itself never lets a user duplicate
{"POST", "/v2/spaces/space1/widgets"},
Expand Down Expand Up @@ -484,6 +485,7 @@ func TestV2Routes(t *testing.T) {
{"DELETE", "/v2/spaces/space1/chats/chat1/messages/msg1"},
{"POST", "/v2/spaces/space1/chats/chat1/messages/msg1/reactions"},
{"POST", "/v2/spaces/space1/chats/chat1/read"},
{"POST", "/v2/spaces/space1/chats/chat1/status"},
{"GET", "/v2/spaces/space1/widgets"},
{"POST", "/v2/spaces/space1/widgets"},
{"PATCH", "/v2/spaces/space1/widgets/w1"},
Expand Down
10 changes: 10 additions & 0 deletions core/api/v2/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -226,6 +226,16 @@ read: no `Idempotency-Key`, `dry_run` ignored.

## Chats

- `POST …/chats/{id}/status` publishes ephemeral activity. An empty body or
`{}` signals typing; the server omits empty `text` so clients can localize
the label. Agents can send `{"text":"Searching documentation","data":{"tool_call":"web_search"}}`.
`data` accepts any JSON value. The encoded payload is limited to 65,508 bytes.
The pubsub topic is `status/<chat_id>`; `Event.Pubsub.Message` carries the
verified sender identity. These updates are not stored messages or SSE
message events. Refresh while active (two seconds recommended), stop when
finished, and let receivers expire status (ten seconds recommended).
`?dry_run=true` does not publish; use a fresh `Idempotency-Key` per refresh.

- `GET …/chats/{id}/messages` returns `{messages, state, message_count,
has_more, next_before?, next_after?}`. `state` carries `unread_messages`,
`unread_mentions`, `last_state_id` — so "anything new?" is a `?limit=1`
Expand Down
1 change: 1 addition & 0 deletions core/api/v2/authz.go
Original file line number Diff line number Diff line change
Expand Up @@ -165,6 +165,7 @@ var v2RouteAuthz = map[string]RouteAuthz{
routeKey(http.MethodDelete, "/v2/spaces/:space_id/chats/:chat_id/messages/:message_id"): {Verb: RouteVerbWrite},
routeKey(http.MethodPost, "/v2/spaces/:space_id/chats/:chat_id/messages/:message_id/reactions"): {Verb: RouteVerbWrite},
routeKey(http.MethodPost, "/v2/spaces/:space_id/chats/:chat_id/read"): {Verb: RouteVerbWrite},
routeKey(http.MethodPost, "/v2/spaces/:space_id/chats/:chat_id/status"): {Verb: RouteVerbWrite},
routeKey(http.MethodPost, "/v2/spaces/:space_id/objects/:object_id/discussion"): {Verb: RouteVerbWrite},
// pairing, registered outside the authenticated group by server.registerAuthRoutes
routeKey(http.MethodPost, "/v2/auth/challenges"): {Verb: RouteVerbWrite, Global: GlobalAuthExempt},
Expand Down
11 changes: 6 additions & 5 deletions core/api/v2/authz_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -254,11 +254,12 @@ func TestV2RouteAuthzTable(t *testing.T) {

t.Run("the non-obvious verb calls hold", func(t *testing.T) {
want := map[string]RouteVerb{
"POST /v2/validate": RouteVerbRead,
"POST /v2/search": RouteVerbRead,
"POST /v2/spaces/:space_id/search": RouteVerbRead,
"POST /v2/spaces/:space_id/chats/:chat_id/read": RouteVerbWrite,
"POST /v2/spaces": RouteVerbWrite,
"POST /v2/validate": RouteVerbRead,
"POST /v2/search": RouteVerbRead,
"POST /v2/spaces/:space_id/search": RouteVerbRead,
"POST /v2/spaces/:space_id/chats/:chat_id/read": RouteVerbWrite,
"POST /v2/spaces/:space_id/chats/:chat_id/status": RouteVerbWrite,
"POST /v2/spaces": RouteVerbWrite,
}
table := RouteAuthzTable()
for key, verb := range want {
Expand Down
82 changes: 82 additions & 0 deletions core/api/v2/handler/chat_status.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,82 @@
package v2handler

import (
"bytes"
"encoding/json"
"fmt"
"io"
"net/http"

"github.com/gin-gonic/gin"

v2model "github.com/anyproto/anytype-heart/core/api/v2/model"
v2service "github.com/anyproto/anytype-heart/core/api/v2/service"
)

// PublishChatStatusHandler publishes ephemeral chat activity.
//
// @Summary Publish chat status
// @Description Publishes encrypted ephemeral JSON on status/<chat_id>. Empty text is omitted so clients can show localized typing. data accepts any JSON value. The encoded payload is limited to 65,508 bytes. Status is not stored or replayed; receivers expire it. Success acknowledges acceptance, not delivery.
// @Id publish_chat_status
// @Tags Chat
// @Accept json
// @Produce json
// @Param space_id path string true "Space id"
// @Param chat_id path string true "Chat or discussion id"
// @Param status body object false "Optional text and arbitrary JSON data; clients localize the default typing label"
// @Param dry_run query bool false "Validate without publishing"
// @Param Idempotency-Key header string false "Unique key for this update; a reused key replays the response without publishing again"
// @Success 200 {object} v2model.ChatStatusResult "Status accepted, or validated on a dry run"
// @Failure 400 {object} v2model.Error "Invalid status"
// @Failure 404 {object} v2model.Error "Space or chat not found"
// @Failure 413 {object} v2model.Error "Request body or encoded pubsub payload too large"
// @Failure 500 {object} v2model.Error "Publication failed"
// @Security bearerauth
// @Router /v2/spaces/{space_id}/chats/{chat_id}/status [post]
func PublishChatStatusHandler(s *v2service.Service) gin.HandlerFunc {
return func(c *gin.Context) {
req, err := decodeChatStatus(c.Request.Body)
if err != nil {
RespondError(c, err)
return
}
result, err := s.PublishChatStatus(c.Request.Context(), c.Param("space_id"), c.Param("chat_id"), req, isV2DryRun(c))
if err != nil {
RespondError(c, err)
return
}
c.JSON(http.StatusOK, result)
}
}

// This body is optional, unlike stored chat messages. Only its data member is
// arbitrary JSON; the envelope remains strict and must be a single object.
func decodeChatStatus(r io.Reader) (req v2model.ChatStatusRequest, err error) {
body, err := io.ReadAll(io.LimitReader(r, maxChatRequestBody+1))
if err != nil {
return req, v2model.ValidationFailed("read status request body", v2model.Issue{Message: err.Error()})
}
if len(body) > maxChatRequestBody {
return req, v2model.RequestTooLarge(fmt.Sprintf("chat status request body exceeds the %d-byte limit", maxChatRequestBody))
}
body = bytes.TrimSpace(body)
if len(body) == 0 {
return req, nil
}
if body[0] != '{' {
return req, v2model.ValidationFailed("status body must be a JSON object")
}
dec := json.NewDecoder(bytes.NewReader(body))
dec.DisallowUnknownFields()
if err := dec.Decode(&req); err != nil {
issue := v2model.Issue{Message: err.Error(), Hint: "send optional text and data; omit text for a localized typing label"}
if field, ok := unknownFieldName(err); ok {
issue.Path = "/" + field
}
return req, v2model.ValidationFailed("invalid status request body", issue)
}
if err := dec.Decode(new(any)); err != io.EOF {
return req, v2model.ValidationFailed("status body must contain exactly one JSON object")
}
return req, nil
}
Loading
Loading