1use objectstore_service::backend::common::{ExpiryUpdate, SetExpiryResponse};
2use objectstore_service::id::{ObjectContext, ObjectId};
3use objectstore_service::multipart::{
4 AbortMultipartResponse, CompleteMultipartResponse, CompletedPart, InitiateMultipartResponse,
5 ListPartsResponse, PartNumber, UploadId, UploadPartResponse,
6};
7use objectstore_service::service::{
8 CreateUploadSessionResponse, DeleteResponse, GetResponse, InsertResponse, MetadataResponse,
9};
10
11use objectstore_service::{ClientStream, StorageService};
12use objectstore_types::auth::Permission;
13use objectstore_types::metadata::Metadata;
14use objectstore_types::range::ByteRange;
15use objectstore_types::resumable::{SessionToken, UploadProgress};
16use objectstore_types::time::Timestamp;
17
18use crate::auth::AuthContext;
19use crate::endpoints::common::ApiResult;
20
21#[derive(Debug)]
47pub struct AuthAwareService {
48 service: StorageService,
49 context: AuthContext,
50 enforce: bool,
51}
52
53impl AuthAwareService {
54 pub fn new(service: StorageService, context: AuthContext, enforce: bool) -> Self {
62 Self {
63 service,
64 context,
65 enforce,
66 }
67 }
68
69 pub fn check_permission(&self, perm: Permission, context: &ObjectContext) -> ApiResult<()> {
73 if let Err(error) = self.context.assert_authorized(perm, context) {
74 sentry::with_scope(
75 |s| s.set_tag("perm", perm.to_string()),
76 || error.log(!self.enforce),
77 );
78
79 if self.enforce {
80 return Err(error.into());
81 }
82 }
83
84 Ok(())
85 }
86
87 pub async fn insert_object(
89 &self,
90 context: ObjectContext,
91 key: Option<String>,
92 metadata: Metadata,
93 stream: ClientStream,
94 access_time: Timestamp,
95 ) -> ApiResult<InsertResponse> {
96 self.check_permission(Permission::ObjectWrite, &context)?;
97 Ok(self
98 .service
99 .insert_object(context, key, metadata, stream, access_time)
100 .await?)
101 }
102
103 pub async fn get_metadata(
105 &self,
106 id: ObjectId,
107 access_time: Timestamp,
108 ) -> ApiResult<MetadataResponse> {
109 self.check_permission(Permission::ObjectRead, id.context())?;
110 Ok(self.service.get_metadata(id, access_time).await?)
111 }
112
113 pub async fn set_expiry(
115 &self,
116 id: ObjectId,
117 target: ExpiryUpdate,
118 access_time: Timestamp,
119 ) -> ApiResult<SetExpiryResponse> {
120 self.check_permission(Permission::ObjectWrite, id.context())?;
121 Ok(self.service.set_expiry(id, target, access_time).await?)
122 }
123
124 pub async fn get_object(
126 &self,
127 id: ObjectId,
128 access_time: Timestamp,
129 range: Option<ByteRange>,
130 ) -> ApiResult<GetResponse> {
131 self.check_permission(Permission::ObjectRead, id.context())?;
132 Ok(self.service.get_object(id, access_time, range).await?)
133 }
134
135 pub async fn delete_object(
137 &self,
138 id: ObjectId,
139 access_time: Timestamp,
140 ) -> ApiResult<DeleteResponse> {
141 self.check_permission(Permission::ObjectDelete, id.context())?;
142 Ok(self.service.delete_object(id, access_time).await?)
143 }
144
145 pub async fn initiate_multipart(
149 &self,
150 id: ObjectId,
151 metadata: Metadata,
152 ) -> ApiResult<InitiateMultipartResponse> {
153 self.check_permission(Permission::ObjectWrite, id.context())?;
154 Ok(self.service.initiate_multipart(id, metadata).await?)
155 }
156
157 pub async fn upload_part(
159 &self,
160 id: ObjectId,
161 upload_id: UploadId,
162 part_number: PartNumber,
163 content_length: u64,
164 content_md5: Option<String>,
165 body: ClientStream,
166 ) -> ApiResult<UploadPartResponse> {
167 self.check_permission(Permission::ObjectWrite, id.context())?;
168 Ok(self
169 .service
170 .upload_part(
171 id,
172 upload_id,
173 part_number,
174 content_length,
175 content_md5,
176 body,
177 )
178 .await?)
179 }
180
181 pub async fn list_parts(
183 &self,
184 id: ObjectId,
185 upload_id: UploadId,
186 max_parts: Option<u32>,
187 part_number_marker: Option<PartNumber>,
188 ) -> ApiResult<ListPartsResponse> {
189 self.check_permission(Permission::ObjectWrite, id.context())?;
190 Ok(self
191 .service
192 .list_parts(id, upload_id, max_parts, part_number_marker)
193 .await?)
194 }
195
196 pub async fn abort_multipart(
198 &self,
199 id: ObjectId,
200 upload_id: UploadId,
201 ) -> ApiResult<AbortMultipartResponse> {
202 self.check_permission(Permission::ObjectWrite, id.context())?;
203 Ok(self.service.abort_multipart(id, upload_id).await?)
204 }
205
206 pub async fn complete_multipart(
208 &self,
209 id: ObjectId,
210 upload_id: UploadId,
211 parts: Vec<CompletedPart>,
212 access_time: Timestamp,
213 ) -> ApiResult<CompleteMultipartResponse> {
214 self.check_permission(Permission::ObjectWrite, id.context())?;
215 Ok(self
216 .service
217 .complete_multipart(id, upload_id, parts, access_time)
218 .await?)
219 }
220
221 pub async fn create_upload_session(
225 &self,
226 id: ObjectId,
227 metadata: Metadata,
228 total_length: u64,
229 ) -> ApiResult<Option<CreateUploadSessionResponse>> {
230 self.check_permission(Permission::ObjectWrite, id.context())?;
231 Ok(self
232 .service
233 .create_upload_session(id, metadata, total_length)
234 .await?)
235 }
236
237 pub async fn put_chunk(
239 &self,
240 id: ObjectId,
241 token: SessionToken,
242 offset: u64,
243 content_length: u64,
244 body: ClientStream,
245 ) -> ApiResult<UploadProgress> {
246 self.check_permission(Permission::ObjectWrite, id.context())?;
247 Ok(self
248 .service
249 .put_chunk(id, token, offset, content_length, body)
250 .await?)
251 }
252
253 pub async fn upload_offset(
255 &self,
256 id: ObjectId,
257 token: SessionToken,
258 ) -> ApiResult<UploadProgress> {
259 self.check_permission(Permission::ObjectWrite, id.context())?;
261 Ok(self.service.upload_offset(id, token).await?)
262 }
263
264 pub async fn cancel_upload(&self, id: ObjectId, token: SessionToken) -> ApiResult<()> {
266 self.check_permission(Permission::ObjectWrite, id.context())?;
268 Ok(self.service.cancel_upload(id, token).await?)
269 }
270}