# Chat Core — Starlark Backend (v0.1.0) # # Library package providing conversations, messages, participants, # and read cursors. Consumable via lib.require("chat-core"). # # Entry points: # on_request(req) → REST API routes # Exported globals → create, send, history, add_participant, # remove_participant, mark_read # # Modules: db, json, realtime # ═══════════════════════════════════════════════ # Helpers # ═══════════════════════════════════════════════ def _resp(status, data): return {"status": status, "body": json.encode(data), "headers": {"Content-Type": "application/json"}} def _str(v): if v == None: return "" return str(v) def _int(v): if v == None: return 0 s = str(v) if not s: return 0 return int(s) def _is_participant(conversation_id, user_id): """Check if user is a participant in the conversation.""" rows = db.query("participants", filters={"conversation_id": conversation_id, "participant_id": user_id}, limit=1) return len(rows or []) > 0 def _get_participant(conversation_id, user_id): """Get participant record or None.""" rows = db.query("participants", filters={"conversation_id": conversation_id, "participant_id": user_id}, limit=1) if rows and len(rows) > 0: return rows[0] return None def _is_admin(conversation_id, user_id): """Check if user has admin role in the conversation.""" p = _get_participant(conversation_id, user_id) if p: return _str(p.get("role", "")) == "admin" return False def _get_conversation(conversation_id): """Get conversation by ID or None.""" rows = db.query("conversations", filters={"id": conversation_id}, limit=1) if rows and len(rows) > 0: return rows[0] return None def _now(): """Current timestamp placeholder — db auto-populates created_at.""" return "" # ═══════════════════════════════════════════════ # Exported API (lib.require("chat-core")) # ═══════════════════════════════════════════════ def create(title, type="group", participants=None, creator_id="", creator_display_name=""): """Create a conversation and add initial participants. Args: title: conversation title type: "direct" or "group" (default "group") participants: list of dicts with {id, type?, display_name?, role?} creator_id: user ID of the creator (added as admin) creator_display_name: display name of the creator Returns: dict with conversation data """ if type not in ("direct", "group"): type = "group" conv = db.insert("conversations", { "title": _str(title), "type": type, "created_by": _str(creator_id), "updated_at": "", }) cid = conv["id"] # Add creator as admin participant if creator_id: db.insert("participants", { "conversation_id": cid, "participant_id": _str(creator_id), "participant_type": "user", "display_name": _str(creator_display_name), "role": "admin", "joined_at": "", }) # Add initial participants for p in (participants or []): pid = _str(p.get("id", "")) if not pid or pid == _str(creator_id): continue db.insert("participants", { "conversation_id": cid, "participant_id": pid, "participant_type": _str(p.get("type", "user")), "display_name": _str(p.get("display_name", "")), "role": _str(p.get("role", "member")), "joined_at": "", }) return conv def send(conversation_id, participant_id, content, content_type="text"): """Send a message to a conversation. Args: conversation_id: target conversation participant_id: sender ID content: message content content_type: "text", "system", or "file" (default "text") Returns: dict with message data """ if content_type not in ("text", "system", "file"): content_type = "text" msg = db.insert("messages", { "conversation_id": _str(conversation_id), "participant_id": _str(participant_id), "content": _str(content), "content_type": content_type, "edited_at": "", }) # Update conversation timestamp db.update("conversations", _str(conversation_id), {"updated_at": msg.get("created_at", "")}) # Publish realtime event realtime.publish( "conversation:" + _str(conversation_id), "message", { "id": msg.get("id", ""), "conversation_id": _str(conversation_id), "participant_id": _str(participant_id), "content": _str(content), "content_type": content_type, "created_at": msg.get("created_at", ""), }, ) return msg def history(conversation_id, limit=50, cursor=""): """Get paginated message history for a conversation. Args: conversation_id: target conversation limit: max messages to return (1-100, default 50) cursor: created_at value of last message from previous page Returns: dict with {messages, has_more, next_cursor} """ lim = _int(limit) if lim < 1 or lim > 100: lim = 50 filters = {"conversation_id": _str(conversation_id)} before = {} if cursor: before = {"created_at": _str(cursor)} rows = db.query("messages", filters=filters, order="-created_at", limit=lim + 1, before=before) rows = rows or [] has_more = len(rows) > lim if has_more: rows = rows[:lim] next_cursor = "" if has_more and rows: next_cursor = _str(rows[-1].get("created_at", "")) return {"messages": rows, "has_more": has_more, "next_cursor": next_cursor} def add_participant(conversation_id, participant_id, participant_type="user", display_name="", role="member"): """Add a participant to a conversation. Returns: dict with participant data """ cid = _str(conversation_id) pid = _str(participant_id) # Check if already a participant existing = db.query("participants", filters={"conversation_id": cid, "participant_id": pid}, limit=1) if existing and len(existing) > 0: return existing[0] p = db.insert("participants", { "conversation_id": cid, "participant_id": pid, "participant_type": _str(participant_type), "display_name": _str(display_name), "role": _str(role), "joined_at": "", }) # Publish realtime event realtime.publish( "conversation:" + cid, "participant.added", { "conversation_id": cid, "participant_id": pid, "participant_type": _str(participant_type), "display_name": _str(display_name), "role": _str(role), }, ) # Insert system message send(cid, pid, "joined the conversation", "system") return p def remove_participant(conversation_id, participant_id): """Remove a participant from a conversation. Returns: True on success """ cid = _str(conversation_id) pid = _str(participant_id) # Delete participant record rows = db.query("participants", filters={"conversation_id": cid, "participant_id": pid}, limit=1) for r in (rows or []): db.delete("participants", r["id"]) # Delete read cursor cursors = db.query("read_cursors", filters={"conversation_id": cid, "participant_id": pid}, limit=1) for c in (cursors or []): db.delete("read_cursors", c["id"]) # Publish realtime event realtime.publish( "conversation:" + cid, "participant.removed", {"conversation_id": cid, "participant_id": pid}, ) return True def mark_read(conversation_id, participant_id, last_read_message_id): """Update the read cursor for a participant. Returns: True on success """ cid = _str(conversation_id) pid = _str(participant_id) mid = _str(last_read_message_id) # Delete existing cursor (upsert via delete+insert) existing = db.query("read_cursors", filters={"conversation_id": cid, "participant_id": pid}, limit=1) for r in (existing or []): db.delete("read_cursors", r["id"]) db.insert("read_cursors", { "conversation_id": cid, "participant_id": pid, "last_read_message_id": mid, }) return True # ═══════════════════════════════════════════════ # REST API dispatcher # ═══════════════════════════════════════════════ def on_request(req): path = req["path"] method = req["method"] user_id = req.get("user_id", "") # ── Conversations ────────────────────────── # GET /conversations — list user's conversations if method == "GET" and path == "/conversations": return _handle_list_conversations(req, user_id) # POST /conversations — create if method == "POST" and path == "/conversations": return _handle_create_conversation(req, user_id) # GET /unread — unread counts if method == "GET" and path == "/unread": return _handle_unread(user_id) # GET /conversations/:id if method == "GET" and path.startswith("/conversations/"): cid = path[len("/conversations/"):] return _handle_get_conversation(cid, user_id) # PUT /conversations/:id if method == "PUT" and path.startswith("/conversations/"): cid = path[len("/conversations/"):] return _handle_update_conversation(cid, req, user_id) # DELETE /conversations/:id if method == "DELETE" and path.startswith("/conversations/"): cid = path[len("/conversations/"):] return _handle_delete_conversation(cid, user_id) # ── Messages ─────────────────────────────── # Routes: /messages/:conversation_id or /messages/:conversation_id/:message_id if path.startswith("/messages/"): remainder = path[len("/messages/"):] parts = remainder.split("/") cid = parts[0] mid = parts[1] if len(parts) > 1 else "" if method == "GET" and not mid: return _handle_list_messages(cid, req, user_id) if method == "POST" and not mid: return _handle_send_message(cid, req, user_id) if method == "PUT" and mid: return _handle_edit_message(cid, mid, req, user_id) if method == "DELETE" and mid: return _handle_delete_message(cid, mid, user_id) # ── Participants ─────────────────────────── if path.startswith("/participants/"): remainder = path[len("/participants/"):] parts = remainder.split("/") cid = parts[0] pid = parts[1] if len(parts) > 1 else "" if method == "GET" and not pid: return _handle_list_participants(cid, user_id) if method == "POST" and not pid: return _handle_add_participant(cid, req, user_id) if method == "DELETE" and pid: return _handle_remove_participant(cid, pid, user_id) # ── Read cursors ────────────────────────── if method == "POST" and path.startswith("/read/"): cid = path[len("/read/"):] return _handle_mark_read(cid, req, user_id) return _resp(404, {"error": "not found"}) # ═══════════════════════════════════════════════ # Conversation handlers # ═══════════════════════════════════════════════ def _handle_list_conversations(req, user_id): """List conversations the user is a participant in.""" # Get all conversation IDs for this user my_parts = db.query("participants", filters={"participant_id": user_id}, limit=500) if not my_parts: return _resp(200, {"data": []}) cids = [] for p in my_parts: cids.append(_str(p.get("conversation_id", ""))) # Fetch conversations and enrich with last message items = [] for cid in cids: conv = _get_conversation(cid) if not conv: continue # Get last message last_msgs = db.query("messages", filters={"conversation_id": cid}, order="-created_at", limit=1) last_msg = None if last_msgs and len(last_msgs) > 0: last_msg = { "id": last_msgs[0].get("id", ""), "content": _str(last_msgs[0].get("content", "")), "participant_id": _str(last_msgs[0].get("participant_id", "")), "content_type": _str(last_msgs[0].get("content_type", "")), "created_at": _str(last_msgs[0].get("created_at", "")), } # Get participant count parts = db.query("participants", filters={"conversation_id": cid}, limit=500) part_count = len(parts or []) items.append({ "id": conv.get("id", ""), "title": conv.get("title", ""), "type": conv.get("type", ""), "created_by": conv.get("created_by", ""), "updated_at": conv.get("updated_at", ""), "created_at": conv.get("created_at", ""), "last_message": last_msg, "participant_count": part_count, }) # Sort by updated_at descending (most recent first) items = sorted(items, key=lambda x: x.get("updated_at", "") or x.get("created_at", ""), reverse=True) return _resp(200, {"data": items}) def _handle_create_conversation(req, user_id): """Create a new conversation.""" body = json.decode(req.get("body", "{}")) title = _str(body.get("title", "")) conv_type = _str(body.get("type", "group")) participants_list = body.get("participants", []) creator_name = _str(body.get("creator_display_name", "")) conv = create(title, conv_type, participants_list, user_id, creator_name) return _resp(201, conv) def _handle_get_conversation(cid, user_id): """Get conversation detail with participants.""" if not _is_participant(cid, user_id): return _resp(403, {"error": "not a participant"}) conv = _get_conversation(cid) if not conv: return _resp(404, {"error": "conversation not found"}) parts = db.query("participants", filters={"conversation_id": cid}, limit=500) return _resp(200, { "id": conv.get("id", ""), "title": conv.get("title", ""), "type": conv.get("type", ""), "created_by": conv.get("created_by", ""), "updated_at": conv.get("updated_at", ""), "created_at": conv.get("created_at", ""), "participants": parts or [], }) def _handle_update_conversation(cid, req, user_id): """Update conversation (admin only).""" if not _is_admin(cid, user_id): return _resp(403, {"error": "admin only"}) body = json.decode(req.get("body", "{}")) patch = {} title = body.get("title", None) if title != None: patch["title"] = _str(title) if patch: db.update("conversations", cid, patch) conv = _get_conversation(cid) return _resp(200, conv) def _handle_delete_conversation(cid, user_id): """Delete conversation and all associated data (admin only).""" if not _is_admin(cid, user_id): return _resp(403, {"error": "admin only"}) # Cascade delete: messages, participants, read_cursors, then conversation msgs = db.query("messages", filters={"conversation_id": cid}, limit=5000) for m in (msgs or []): db.delete("messages", m["id"]) parts = db.query("participants", filters={"conversation_id": cid}, limit=500) for p in (parts or []): db.delete("participants", p["id"]) cursors = db.query("read_cursors", filters={"conversation_id": cid}, limit=500) for c in (cursors or []): db.delete("read_cursors", c["id"]) db.delete("conversations", cid) return _resp(200, {"ok": True}) # ═══════════════════════════════════════════════ # Message handlers # ═══════════════════════════════════════════════ def _handle_list_messages(cid, req, user_id): """Paginated message history.""" if not _is_participant(cid, user_id): return _resp(403, {"error": "not a participant"}) q = req.get("query", {}) limit_val = _int(q.get("limit", "50")) cursor = _str(q.get("cursor", "")) result = history(cid, limit_val, cursor) return _resp(200, result) def _handle_send_message(cid, req, user_id): """Send a message.""" if not _is_participant(cid, user_id): return _resp(403, {"error": "not a participant"}) body = json.decode(req.get("body", "{}")) content = _str(body.get("content", "")) content_type = _str(body.get("content_type", "text")) if not content: return _resp(400, {"error": "content required"}) msg = send(cid, user_id, content, content_type) return _resp(201, msg) def _handle_edit_message(cid, mid, req, user_id): """Edit a message (author only).""" if not _is_participant(cid, user_id): return _resp(403, {"error": "not a participant"}) # Verify ownership msgs = db.query("messages", filters={"id": mid}, limit=1) if not msgs: return _resp(404, {"error": "message not found"}) msg = msgs[0] if _str(msg.get("participant_id", "")) != user_id: return _resp(403, {"error": "can only edit own messages"}) body = json.decode(req.get("body", "{}")) content = _str(body.get("content", "")) if not content: return _resp(400, {"error": "content required"}) # Use created_at of the message as the edited_at marker db.update("messages", mid, {"content": content, "edited_at": msg.get("created_at", "")}) # Publish edit event realtime.publish( "conversation:" + cid, "message.edited", {"id": mid, "conversation_id": cid, "content": content, "edited_by": user_id}, ) updated = db.query("messages", filters={"id": mid}, limit=1) return _resp(200, updated[0] if updated else {}) def _handle_delete_message(cid, mid, user_id): """Delete a message (author or admin).""" if not _is_participant(cid, user_id): return _resp(403, {"error": "not a participant"}) msgs = db.query("messages", filters={"id": mid}, limit=1) if not msgs: return _resp(404, {"error": "message not found"}) msg = msgs[0] is_author = _str(msg.get("participant_id", "")) == user_id is_conv_admin = _is_admin(cid, user_id) if not is_author and not is_conv_admin: return _resp(403, {"error": "can only delete own messages or be admin"}) db.delete("messages", mid) # Publish delete event realtime.publish( "conversation:" + cid, "message.deleted", {"id": mid, "conversation_id": cid, "deleted_by": user_id}, ) return _resp(200, {"ok": True}) # ═══════════════════════════════════════════════ # Participant handlers # ═══════════════════════════════════════════════ def _handle_list_participants(cid, user_id): """List participants in a conversation.""" if not _is_participant(cid, user_id): return _resp(403, {"error": "not a participant"}) parts = db.query("participants", filters={"conversation_id": cid}, limit=500) return _resp(200, {"data": parts or []}) def _handle_add_participant(cid, req, user_id): """Add a participant (admin only).""" if not _is_admin(cid, user_id): return _resp(403, {"error": "admin only"}) body = json.decode(req.get("body", "{}")) pid = _str(body.get("participant_id", "")) if not pid: return _resp(400, {"error": "participant_id required"}) ptype = _str(body.get("participant_type", "user")) display_name = _str(body.get("display_name", "")) role = _str(body.get("role", "member")) p = add_participant(cid, pid, ptype, display_name, role) return _resp(201, p) def _handle_remove_participant(cid, pid, user_id): """Remove a participant (admin only, or self-remove).""" is_self = pid == user_id if not is_self and not _is_admin(cid, user_id): return _resp(403, {"error": "admin only (or remove yourself)"}) remove_participant(cid, pid) return _resp(200, {"ok": True}) # ═══════════════════════════════════════════════ # Read cursor handlers # ═══════════════════════════════════════════════ def _handle_mark_read(cid, req, user_id): """Mark conversation as read up to a message.""" if not _is_participant(cid, user_id): return _resp(403, {"error": "not a participant"}) body = json.decode(req.get("body", "{}")) mid = _str(body.get("last_read_message_id", "")) if not mid: return _resp(400, {"error": "last_read_message_id required"}) mark_read(cid, user_id, mid) return _resp(200, {"ok": True}) def _handle_unread(user_id): """Get unread counts for all user's conversations.""" my_parts = db.query("participants", filters={"participant_id": user_id}, limit=500) if not my_parts: return _resp(200, {"data": {}}) counts = {} for p in my_parts: cid = _str(p.get("conversation_id", "")) if not cid: continue # Get read cursor cursors = db.query("read_cursors", filters={"conversation_id": cid, "participant_id": user_id}, limit=1) if cursors and len(cursors) > 0: last_read_id = _str(cursors[0].get("last_read_message_id", "")) if last_read_id: # Get the created_at of the last read message last_read_msgs = db.query("messages", filters={"id": last_read_id}, limit=1) if last_read_msgs and len(last_read_msgs) > 0: last_read_at = _str(last_read_msgs[0].get("created_at", "")) # Count messages after the read cursor unread = db.query("messages", filters={"conversation_id": cid}, after={"created_at": last_read_at}, limit=1000) counts[cid] = len(unread or []) else: # Last read message was deleted — count all all_msgs = db.query("messages", filters={"conversation_id": cid}, limit=1000) counts[cid] = len(all_msgs or []) else: # No cursor value — all messages are unread all_msgs = db.query("messages", filters={"conversation_id": cid}, limit=1000) counts[cid] = len(all_msgs or []) else: # No cursor at all — all messages are unread all_msgs = db.query("messages", filters={"conversation_id": cid}, limit=1000) counts[cid] = len(all_msgs or []) return _resp(200, {"data": counts})